Skip to content

Persist full job completion batches concurrently - #1434

Draft
bgentry wants to merge 2 commits into
masterfrom
bg/parallel-job-completion
Draft

bgentry wants to merge 2 commits into
masterfrom
bg/parallel-job-completion

Conversation

@bgentry

@bgentry bgentry commented Oct 5, 2026

Copy link
Copy Markdown
Contributor

DRAFT: not ready for review. Opened to evaluate on its own merits and schedule separately.

The batch completer persists finished jobs with one JobSetStateIfRunningMany query at a time. That works well under light load, but under sustained heavy load it becomes the bottleneck. While one 5,000-job batch is being written, workers keep finishing jobs and a second full batch piles up behind it. When the backlog reaches its limit, completions start waiting on the completer, so workers sit idle even though the database could easily run another query.

Here's an example: a client works down a large queue of short jobs. Each completion query takes a few milliseconds, and in that time workers finish thousands more jobs. The next batch is already full by the time the current query returns, so the completer is always behind, and job throughput is capped by the latency of one query at a time.

With this change, the completer can run up to two completion queries at once, each still capped at 5,000 jobs. The second query starts only after its own backlog has filled a batch, so light and sparse traffic keeps coalescing exactly as it does today and doesn't use an extra connection. Concurrency is opt-in through two new optional interfaces, riverdriver.ExecutorJobCompletionConcurrency and riverpilot.PilotJobCompletionConcurrency. The pgx and database/sql pool executors and the standard pilot report two. Transaction executors report one, since one transaction can't run statements concurrently. SQLite and any pilot that doesn't implement the interface stay serial.

Running two batches at once introduces a hazard. Two batches could hold different executions of the same job, and PostgreSQL might lock their rows out of order, which would let an older result overwrite a newer attempt. To prevent that, the completer tracks in-flight and retrying job IDs. A later result for the same job waits until the earlier batch finishes, while batches for different jobs still persist concurrently. On shutdown, the completer waits for active queries and then makes one bounded final flush attempt, as before.

There are two user-visible side effects, both noted in the changelog. Subscribers may get completion events for different jobs in a different order than before, though events for the same job stay in order. And the completer may briefly use one extra database connection.

I compared river bench on master (6cdaf38) and on this branch, with runs alternating between the two builds against the same local database:

river bench -n 1000000 (5 runs each) master this branch
median jobs/sec 58,916 68,416 (+16%)
range 43,542 – 62,508 65,351 – 69,264
river bench --duration 30s (3 runs each) master this branch
mean jobs/sec 48,312 49,450 (+2%)

The -n mode inserts all the jobs up front and then measures how fast they're worked down, which is the backlog case this change targets. The --duration mode inserts and works at the same time, and there the inserter limits throughput, so the gain is small. These numbers come from one laptop (Apple Silicon M4, 14 cores, Homebrew PostgreSQL 18.6 on localhost, Go 1.27.1). Absolute numbers will differ on other machines and on a networked database, where query latency is higher and overlapping queries may help more.

This also brings Go's completion throughput back in line with River's Rust implementation, which persists completion batches the same way.

The first commit adds the capability interfaces and a driver test that runs the advertised number of disjoint completion batches concurrently through a pool executor. The second commit changes the completer, with tests for capability negotiation, a sparse second batch that stays queued until it's full, and a job whose results are forced to interleave across concurrent batches.

@bgentry
bgentry force-pushed the bg/parallel-job-completion branch 3 times, most recently from b67f172 to c744f1d Compare October 6, 2026 18:49
Add optional `ExecutorJobCompletionConcurrency` and
`PilotJobCompletionConcurrency` interfaces through which an executor and
a pilot can declare how many `JobSetStateIfRunningMany` calls they can
safely run at once. Implementations that don't implement them are
treated as serial.

The pgx and `database/sql` pool executors and the standard pilot report
two. Their transaction and subtransaction executors report one, since a
single transaction can't run statements concurrently. SQLite and other
drivers don't opt in.

A new driver test runs the advertised number of disjoint completion
batches concurrently through a pool executor and checks that a
transaction never claims more than one.
The batch completer persists one `JobSetStateIfRunningMany` call at a
time. Under sustained load a second full backlog builds while the first
query is still in flight, so producers wait on the backlog and workers
sit idle even though the database has spare capacity.

When both the executor and the pilot declare completion concurrency,
run up to two completion queries at once, each capped at the existing
5,000-job batch size. A second query starts only once its own backlog
is full, so sparse traffic keeps the existing coalescing and connection
use. SQLite, transactions, and pilots that don't opt in stay serial.

Concurrent batches could otherwise hold different executions of the
same job, and PostgreSQL may lock their rows out of submission order,
letting an older result overwrite a newer attempt. The completer tracks
in-flight and retrying job IDs so a later result for the same job waits
behind the earlier batch while independent batches still persist
concurrently. Shutdown waits for active queries and makes one bounded
final flush attempt.

Subscribers may see completion events for different jobs in a different
order than before, and the completer may briefly use one more database
connection.

Tests block the first query and check that a sparse second batch stays
queued through a flush interval and starts as soon as it reaches the
batch threshold, cover capability negotiation, and force a duplicate job
ID to interleave across batches.
@bgentry
bgentry force-pushed the bg/parallel-job-completion branch from c744f1d to adba3de Compare October 7, 2026 13:54
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant